-----------------------------------------------------------------------------------------
---  Dstream : reduceByKey
-----------------------------------------------------------------------------------------
# TERMINAL 1:
-------------
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._ // not necessary since Spark 1.3

val conf = new SparkConf().setMaster("local[2]").setAppName("NetworkWordCount").set("spark.driver.allowMultipleContexts","true")
val ssc = new StreamingContext(conf, Seconds(3))

val lines = ssc.socketTextStream("localhost", 9988)

val words = lines.flatMap(_.split(" "))

import org.apache.spark.streaming.StreamingContext._ // not necessary since Spark 1.3

// Count each word in each batch
val pairs = words.map(word => (word, 1))
val wordCounts = pairs.reduceByKey(_ + _)

wordCounts.print()

// Start the computation
ssc.start()             

# TERMINAL 2:
-------------
// permet d'crire les mots qui seront ensuite traits
nc -lk 9988 





-----------------------------------------------------------------------------------------
--- file foreachRDD
-----------------------------------------------------------------------------------------
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._

val conf = new SparkConf().setAppName("File Count").setMaster("local[2]").set("spark.driver.allowMultipleContexts","true")

val sc = new SparkContext(conf)
val ssc = new StreamingContext(sc, Seconds(1))

val file = ssc.textFileStream("/opt/spark/examples/src/main/resources/file1")
file.foreachRDD(t=> {

val test = t.flatMap(line => line.split(" ")).map(word => (word, 1)).reduceByKey(_ + _)

test.saveAsTextFile("/opt/spark/examples/src/main/resources/file2")

                    })
ssc.start()


-----------------------------------------------------------------------------------------
--- Dstream : Sliding Window 
-----------------------------------------------------------------------------------------
# TERMINAL 1:
-------------

import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.StreamingContext._ // not necessary since Spark 1.3

val conf = new SparkConf().setMaster("local[2]").setAppName("NetworkWordCount").set("spark.driver.allowMultipleContexts","true")
val ssc = new StreamingContext(conf, Seconds(3))

val mystream = ssc.socketTextStream("localhost", 9988)

//On applique ce traitement sur l'ensemble des DStream datant de moins d'une
//minutes. Le calcul est refait toute les 4 secondes


var words = mystream.flatMap(line => line.split(" ")).map(x=>(x,1))

val nbMot = words.reduceByKeyAndWindow((x:Int,y:Int)=>x+y,Seconds(9),Seconds(3))

nbMot.print()

ssc.start()



# TERMINAL 2:
-------------
// permet d'crire les mots qui seront ensuite traits
nc -lk 9988